diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/common/ContextVariable.java b/proxy/src/main/java/org/apache/rocketmq/proxy/common/ContextVariable.java index f77ad376f4..fcc6bb02ff 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/common/ContextVariable.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/common/ContextVariable.java @@ -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"; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/common/ProxyExceptionCode.java b/proxy/src/main/java/org/apache/rocketmq/proxy/common/ProxyExceptionCode.java index 50a5ce1353..662df1b900 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/common/ProxyExceptionCode.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/common/ProxyExceptionCode.java @@ -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, } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcProxyException.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcProxyException.java index ed1eaf1197..f89639dc58 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcProxyException.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcProxyException.java @@ -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) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivity.java index 59bbb0ea0b..ca215e6777 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/producer/SendMessageActivity.java @@ -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 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) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/AbstractProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/AbstractProcessor.java index bf47d2d116..6a97754b01 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/AbstractProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/AbstractProcessor.java @@ -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; + } + } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java index b93187380a..2007fa53cd 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java @@ -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 messageExtList, long timeoutMillis) { CompletableFuture> 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); } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/validator/DefaultTopicMessageTypeValidator.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/validator/DefaultTopicMessageTypeValidator.java index a0718c5ced..eaa4144c5f 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/validator/DefaultTopicMessageTypeValidator.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/validator/DefaultTopicMessageTypeValidator.java @@ -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()); } } }